MQ - RocketMQ 5.4 横向和纵向技术全面探究以及它的安装使用介绍

为什么要引入MQ?

使用 MQ 的三大好处

在实际项目中,如果没有消息中间件,系统之间就是通过 RPC(比如 Dubbo 或简单的 HTTP)进行硬连接。这在业务规模变大后,会撞上三堵物理墙。为了撞碎这三堵墙,衍生出了 MQ 这种异步处理技术:

  • 一:异步处理(提升系统吞吐量)
    • 典型的场景:用户在某系统上注册了一个账号。系统需要做三件事:创建用户记录(10ms)、发送激活邮件(200ms)、发送欢迎短信(150ms)。
    • 如果不用 MQ 而使用 RPC 硬连:用户点击“注册”,前端转圈圈,一共要等 $10 + 200 + 150 = 360\text{ms}$ 才能看到注册成功。如果邮件服务器卡了,用户可能要等 5 秒。
    • 用了 MQ 的工业套路:用户点击 “注册”,后台把用户记录往数据库一塞,然后往 RocketMQ 里丢一条名为 “用户已注册” 的消息(耗时 2ms),直接给前端返回 “注册成功”。至于发邮件、发短信的系统,它们自己去向 RocketMQ 订阅这条消息,在后台默默地异步执行。用户的等待时间从 360ms 缩短到 12ms。
  • 二:系统解耦(提升系统的可维护性和可扩展性)
    • 典型的场景:运营系统、数据分析系统(大屏)、风控系统,全都对“用户下单”这个事件感兴趣。
    • 不用 MQ 的硬连:每当新加一个系统,订单系统的开发人员就要加班,在自己的代码里加一行 rpcToNewSystem(),一旦新系统挂了还得考虑怎么重试。
    • 用了 MQ 的工业套路:订单系统只管往 RocketMQ 发一句:“有人下单了,数据在这”。谁感兴趣谁自己去跟 RocketMQ 打招呼(订阅),订单系统一概不管。
  • 三:削峰填谷(保护脆弱的后端)
    • 典型的场景:电商系统突发秒杀活动,或者大促。平时每秒只有 100 个请求(下订单),系统很清闲。秒杀开始那一秒,瞬间涌入 10,000 个请求。
    • 不用 MQ 的硬连:这 10,000 个请求通过 RPC 毫无阻挡地直接砸向底层的 MySQL 数据库。MySQL 并发上限可能就 2,000,瞬间直接被砸死、瘫痪。
    • 用了 MQ 的工业套路:在网关层收到这 10,000 个请求后,不做任何数据库操作,直接做个简单的合法性校验,然后打包成消息全部塞进 RocketMQ(RocketMQ 的写入并发极其恐怖,能轻松抗住十万级)。此时后端真正的订单消费系统,按照自己数据库能承受的极限(比如每秒 1,000 条),慢条斯理地从 RocketMQ 里拉取请求进行处理。这就是“削峰(把突发大洪峰接住)填谷(在随后的时间里平稳消化)


使用 MQ 带来的四个挑战

我们知道,世界上没有任何一种架构设计是十全十美的,任何优势的背后都暗中标好了代价(Trade-off)。如果使用不当,必然带来坏处,以下四点是使用 MQ 必须考虑到的:

  • 【系统可用性降低】
    • 在没有引入 MQ 之前,你的系统只有 “前端 $\rightarrow$ 订单服务 $\rightarrow$ 数据库” 这三个环节。 引入 MQ 后,链路变成了:“前端 $\rightarrow$ 订单服务 $\rightarrow$ MQ (Broker) $\rightarrow$ 库存/核心业务服务 $\rightarrow$ 数据库”。
    • 系统的整体可用性是各个组件可用性的乘积。链路越长,崩盘的概率越高。原本 MQ 没来时,系统活得好好的。现在一旦 MQ 机器突然断电、硬盘坏了或者集群网络抖动,整个上游的订单投递会全部死掉,导致全线崩溃。你不得不花大代价去搞 RocketMQ 的双主双从、DLedger 自动选举高可用集群来补这个窟窿。
  • 【系统复杂度飙升】
    • 使用了异步处理后,原本一条线走到底的简单同步代码被切成了两截。这会引出一系列让人抓狂的分布式难题:
      • 痛点 A:如何保证消息不丢【不丢】?(网络抖动了,Sender 发成功了吗?Broker 真正落盘了吗?Consumer 消费时死机了怎么办?)
      • 痛点 B:如何解决重复消费【不重】?(因为网络超时,生产者以为没发成功,又重发了一次;或者消费者 Ack 丢了,Broker 又推了一次。这就要求下游所有消费端必须都做幂等性设计,比如去重表、分布式锁)。
      • 痛点 C:如何保证消息的顺序【不乱】?(本来是先下订单,后支付。异步化之后,支付消息跑得比下单消息还快,下游系统直接报错)。
  • 【分布式数据一致性问题】
    • 在传统的同步调用中,我们可以在一个方法里加上 @Transactional 本地事务。订单失败了,库存自动回滚,数据是强一致的。
    • 引入 MQ 异步化后:订单系统成功把消息丢给了 MQ,自己高高兴兴提交了事务。结果库存系统在消费这条消息时,由于数据库刚好卡死,扣减库存失败了!此时订单下了,库存没扣,数据脏了。这就是分布式世界最头疼的 “一致性破坏”。你必须引入复杂的分布式事务(如 RocketMQ 的事务消息机制)、或者人工写补偿脚本、对账程序来确保它们能 “最终一致”。
  • 【业务天然失去了“即时同步响应”】
    • 异步意味着 “不可预期的等待”。有些场景是无法接受异步的。比如 “用户刷卡支付”,用户必须在转圈圈的几秒钟内知道自己到底是付款成功了还是余额不足。你不能对用户说:“您先走吧,我们把付款请求塞进 MQ 了,今晚后台慢慢扣,扣失败了我们再电话通知您。”


如何抉择?

所以,我们在实际业务中到底用不用 MQ,完全取决于一笔账:

$$\text{业务痛点(洪峰压死系统 / 响应太慢)} > \text{引入 MQ 带来的代价(高可用运维 + 幂等设计 + 最终一致性处理)}$$

如果大促的洪峰能把你的数据库直接砸瘫痪(痛点致命),那么哪怕 MQ 再复杂、再有坏处,你也必须硬着头皮上,然后用技术去死磕上述的四个坏处。


从MQ协议角度纵览消息队列

纵观 MQ(消息中间件)的演进史,本质上是一部“通用标准化”与“极致性能化”相互博弈的协议进化史。它主要经历了三个标志性阶段:

  • JMS 时代(Java 内部规范,2001年左右):早期企业级应用缺乏统一的异步通信标准。Sun 公司推出 JMS(Java Message Service),它并非网络传输协议,而是类似 JDBC 的 API 规范。它定义了点对点(Queue)与发布订阅(Topic)两大模型,成功统一了 Java 世界的消息访问,催生了 ActiveMQ 等早期产品。但它绑定 Java 语言,且无法对底层网络和存储进行极致优化。
  • AMQP 时代(国际通用标准化,2006年左右):为了打破跨语言、跨厂商的壁垒,金融巨头与科技大厂联合制定了 AMQP(高级消息队列协议)。这是一种真正的二进制网络应用层协议。它引入了 Exchange(交换机)概念,实现了极其精细、灵活的路由拓扑规则,落地产品代表为 RabbitMQ。但 AMQP 为追求通用与复杂路由,牺牲了高并发下吞吐量上限。
  • 自定义私有协议时代(高并发与大数据时代,2011年至今):互联网洪峰(大数据、秒杀交易)带来的物理事实,让 AMQP 的复杂路由沦为性能瓶颈。以 Kafka 和 RocketMQ 为代表的产品彻底抛弃通用标准,基于 TCP 自研高效的二进制私有协议(如魔数+消息长度+Payload)。它们摒弃复杂的交换机计算,将消息统一顺序追加写入磁盘文件,用极致简单的协议头换取了百万级的吞吐量。同时,物联网的爆发也让短小精悍(2字节固定头)的 MQTT 协议成为了 IoT 领域的绝对标准。


当前 MQ 产品的横向对比

在工业界,消息中间件经历了几代演进。目前市面上最常被提起的五大主流 MQ 分别是:ActiveMQ(老一代)、RabbitMQ、Kafka、RocketMQ、Apache Pulsar。为了让你有最感性的认识,我们把它们比作四种不同类型的“交通工具”:

  • RabbitMQ (高灵敏跑车)
    • Erlang 语言编写。历史悠久,在中小企业、传统金融、物联网中应用极广。
    • 极致的灵活性与低延迟。支持非常复杂的路由规则,消息到了它手里能像流水一样快速分发。
    • 致命短板是:吞吐量较低(万级),Erlang 语言小众,极难深入读源码和定制。
    • 它就像高灵敏跑车,变道灵活、路线复杂,但拉不了太重的货物。
  • Kafka (超长重载运煤火车)
    • 领英(LinkedIn)开源,Scala/Java 编写。大数据领域的绝对霸主。
    • 极致的吞吐量,每秒能吞吐百万级的数据,文件存储设计极其暴力。
    • 致命短板是:队列(Partition)多了之后性能会剧烈下滑;不支持复杂的业务过滤,容易丢数据(不适合金融交易)。
    • 它就像是一列超长的运煤火车,沿直线狂奔,一趟拉几百万吨(大数据日志),但掉头极难,掉几块 “煤渣” 它也不太在乎。
  • RocketMQ (工业级高铁)
    • 阿里自研后捐赠给 Apache,纯 Java 编写。为了电商交易而生的王牌。
    • 极致的稳定性与业务可靠性。在抗住超高并发的同时,绝不丢一条消息。
    • 致命短板是:相比 Kafka,纯粹的吞吐量上限略逊一筹;生态上在大数据领域的适配不如 Kafka 广泛。
    • 它就像是工业级高铁,既有火车的恐怖吞吐量,又具备极高的调度精准度,绝不错发消息(金融级安全)。
  • Pulsar(模块化可拆卸的科幻概念车)
    • Pulsar 绝对是一个不容忽视的庞然大物。它的出现,几乎就是为了精准打击 Kafka 和 RocketMQ 在架构设计上的“物理硬伤”。
    • 计算与存储彻底分离、多租户。天生为云原生、大厂基础架构而生。
    • 以上MQ基本都是存算一体的(RocketMQ5.x 除外),但是 Pusar 是存算分离。Broker 只管计算(无状态),底层的存储丢给另一个专门的组件 Apache BookKeeper(强状态)。
    • 如果把前面的三个 MQ 比作传统的交通工具,那么 Pulsar 就是一台 “模块化可拆卸的科幻概念车”。车头(计算)和车厢(存储)是分离开的。车头不够了直接加车头,车厢不够了直接挂车厢。
    • Pulsar 独特性体现在:
      • 无论是 Kafka 还是 RocketMQ,在工业上面临一个巨大的物理痛点:扩容极其痛苦(重平衡 Rebalance)。比如你的订单 Topic 数据量暴增,当前 3 台 Broker 磁盘都快撑爆了。你紧急买了两台新服务器(Broker 4、Broker 5)加进来。由于它们是存算一体的,新机器上空空如也。为了分担压力,你必须把前 3 台机器上的历史数据(几十个 G 甚至几个 T),通过网络硬生生地复制、搬迁到新机器上。代价是搬迁数据会瞬间榨干机房的网络带宽,导致正常的业务请求大量超时雪崩。
      • Pulsar 的创始人看透了这个痛点,决定在架构上动大手术:
        • Pulsar Broker(车头/计算层):它是完全无状态的。只负责接收生产者的 send() 请求,转手就把数据扔给底层。它不存任何数据!
        • Apache BookKeeper(车厢/存储层):专门负责高效存盘。数据是以一个个小数据块(Ledger)分散存在多个存储节点上的。
        • 当 Pulsar 需要扩容时:如果并发请求太高,CPU 撑不住了:直接增加几台 Pulsar Broker 机器。因为它们无状态,一秒钟就启动好了,不需要搬迁任何数据。
        • 如果磁盘满了,数据装不下了:直接增加几台 BookKeeper 存储机器。新来的数据直接写进新机器,历史数据老老实实呆在原地,根本不用动。
    • 既然 Pulsar 架构这么科幻、这么完美,为什么目前国内大厂的大多数核心业务依然选择 RocketMQ? 这里就涉及到了高维度的 Trade-off(工程权衡) 实践事实:
      • 运维成本是吞金兽:
        • RocketMQ 极度轻量。它就一个 NameServer(几百行代码)和一个 Broker。
        • Pulsar 想要跑起来,你得维护:Pulsar 自身集群 + BookKeeper 存储集群 + ZooKeeper(或 Metadata)集群。为了发个消息,你得同时伺候三套分布式集群。
      • 业务特性没有 RocketMQ 贴心:
        • RocketMQ 是阿里大促里一点点抠出来的,它的事务消息、多级延时消息开箱即用,设计得非常贴合国内程序员的业务直觉。
        • Pulsar 早期更侧重于基础架构的大容量和多租户隔离(比如腾讯等大厂用来做统一的消息云平台),在细粒度的分布式事务业务支持上,生态不如 RocketMQ 那么丝滑。


理解 RocketMQ 的几个术语

RocketMQ 里的名词很多,我们用 “寄快递” 的现实生活场景来做一个生动的解释:

  • Topic(主题) :快递的分流大方向(如:广东省、北京市)。它是消息的逻辑分类。比如订单系统的消息叫 TOPIC_ORDER,用户系统的消息叫 TOPIC_USER。
  • Tag(标签) :快递的具体写明类型(如:生鲜、易碎、普通)。它是 Topic 下面的次级分类,用来进一步过滤。比如 TOPIC_ORDER 下面,可以有 Tag_Pay_Success(支付成功)和 Tag_Cancel(订单取消)。
  • Producer(生产者) :发件人。负责把消息(快递)打包并投递出去的业务系统。
  • Consumer(消费者) :收件人。负责接收消息(接收快递)并执行具体业务的系统。
  • Broker(大管家) :邮局的各个网点/中转站。它是绝对的实体核心!负责接收生产者发来的消息,把消息在磁盘上存好,再把消息高高兴兴地递给消费者。
  • NameServer(协调者) :邮局的全国营业网点通信录/路由表。它负责记录哪台 Broker 在哪个 IP,里面有哪些 Topic。Producer 和 Consumer 在干活前,先去 NameServer 查一下:“我想寄广东省的快递,应该去哪个网点(Broker IP)”。


RocketMQ 5.x 系列的云原生化

目前 RocketMQ 5.x 系列已是绝对的主流首选。在 4.x 的基础上,5.x 做了如下重大升级:

  • 冷热分层存储:
    • 热数据就是刚刚产生的消息(如 5 分钟内),这些消息往往会立刻被下游消费者疯狂拉取和处理。由于刚写进磁盘,它们还热乎着,大量残留在操作系统的 PageCache(内核内存)里,消费者去读时,直接命中内存,速度极快。
    • 冷数据是指过去很久的消息(如 3 天前、1 周前)。根据业务要求,为了对账、回溯或数据分析,大厂往往要求消息在 Broker 上至少保留 3 天到 7 天。
    • 随着时间推移,物理痛点暴露了。磁盘每天产生几个 T 的消息,7 天就是几十个 T。Broker 服务器的高性能固态硬盘(SSD)贵得要死,很快就会被塞满。如果此时有某个败家子业务,突然要从 3 天前的旧消息开始 “重新消费”。Broker 没办法,磁头只能去磁盘的角落里把旧文件读出来。这一读,旧数据会瞬间挤占宝贵的 PageCache(内存),把当前正在高频读写的“热数据”给硬生生挤出去(术语叫 PageCache 颠簸)。导致整个 Broker 的写性能瞬间雪崩,秒杀业务直接卡死。
    • 为了解决 “既要保留很久的历史数据,又不能拖垮当前高性能写盘” 的矛盾,RocketMQ 5.x 设计了冷热分层存储架构。它的核心思想是让高性能 SSD 专职伺候热数据,让便宜的云端对象存储(如阿里云 OSS、AWS S3、MinIO)去接盘冷数据。它的自动运转时序是:
      • 热数据留在本地:生产者发来的最新消息,依然高高兴兴地顺序写入 Broker 本地的 CommitLog(高性能 SSD 盘)。
      • 后台异步 “偷运”:Broker 内部有一个常驻的后台线程。它会像清洁工一样,定时扫描本地文件。一旦发现某些消息文件的修改时间超过了设定的阈值(比如超过 1 小时,变成“冷数据”了),它就会在后台默默地、异步地把这些数据打包上传到低成本的对象存储(如 OSS/S3)中。
      • 本地空间释放:一旦冷数据安全到达云端,Broker 就会把本地 SSD 盘上的历史文件删掉,腾出空间迎接新的热数据。
      • 透明消费:当消费者来拉取消息时,客户端不需要关心里面的逻辑。如果它消费的是最新消息,Broker 直接从本地 SSD 盘或内存里丢给它。如果它突然发起历史回溯,去消费 3 天前的消息,Broker 发现本地没有,会自动、透明地去云端对象存储把数据拉回来发给消费者。
      • 分层存储出来后:本地 Broker 只需要买 1TB 的 SSD 硬盘,只要能抗住最近几小时的洪峰就行。剩下的 49TB 全部自动归档到对象存储(OSS)。对象存储的价格只有 SSD 硬盘的 1/10 甚至更低,且天然自带无限扩容和多副本高可用属性。这意味着,5.x 的冷热分层存储,在不牺牲秒杀性能的前提下,直接帮企业砍掉了 80% 以上的存储开销。
  • Prox架构:
    • 简单来说就是 RocketMQ 在客户端和 Broker 之间强行加了一层专属服务员叫做 Proxy。客户端(尤其是跨语言、轻量级的客户端)不再直接和底层的 Broker 建立复杂的连接,也不用管数据存在哪个 Broker。客户端只需要和无状态的 Proxy 建立连接(通常走国际通用的 gRPC 协议),把消息交给 Proxy。Proxy 作为一个大管家,再在后台把消息分发给真正的 Broker。
    • 在传统的 4.x 经典 Remoting 架构中,Broker 是一个典型的 “大胖子”:它既要负责处理客户端的连接、网络解析、权限校验(计算),又要负责把消息顺序写盘(存储)。5.x 通过引入无状态的 Proxy 架构,实现了真正的存算分离。计算层(Proxy)完全无状态,专门负责接客、解析协议、鉴权。存储层(Broker)专注于存盘,它被稳稳地保护在后端,只和 Proxy 打交道,不需要面对错综复杂的客户端外网连接。
    • 但是目前的现实是,国内 90% 以上的大厂(如阿里、美团、滴滴)的历史核心业务和金融交易,因为追求极致的吞吐量和低延迟,依然大面积跑在 Remoting 经典架构上。Remoting 架构中客户端直接和 Brocker 打交道,没有中间层 Proxy 赚差价,它是我们理解 RocketMQ 肉身的关键。
  • 全面拥抱 gRPC 协议:
    • 在 4.x 时代,RocketMQ 采用的是纯自研的二进制私有协议,这直接导致了两个限制:
      • 跨语言极难:别的语言(如 Go、Python、C++)想要生态适配,必须硬着头皮用该语言重写一套极其复杂的客户端通信底层(Remoting 驱动)。
      • 云网关接入难:传统的私有协议很难直接通过云原生中的 Service Mesh(如 Istio)或者标准 Service 来进行七层网络治理和流量调度。
    • 5.x 的拥抱姿态:全面拥抱 gRPC 协议。
      • gRPC 是云原生标准基金会(CNCF)的官方大杀器,天然基于 HTTP/2 传输,并且天生支持各种语言客户端的自动生成。
      • 现在,无论是写 Go 的云原生微服务,还是写 Python 的 AI 数据处理流,都可以直接用一行标准 gRPC 代码无缝连接 RocketMQ 5.x,整个框架在云端网络拓扑里变得极为透明和通用。
  • 全方位接入云原生三大可观测性标准:
    • 云原生架构(如 K8s 体系)有一个非常核心的铁律:任何不能被监控、被度量、被追踪的组件,都是云上的定时炸弹。
    • 4.x 的局限:想要看 RocketMQ 的健康指标,你必须去装它专门的配套项目 RocketMQ Console,或者自己去解析它的原生 JMX 监控,很难和全公司的监控大盘(如 Prometheus)融合。
    • 5.x 的拥抱姿态:全方位接入云原生三大可观测性标准。
      • Metrics(指标):原生内置了标准的 Prometheus 监控指标打点,运维人员可以直接在 Grafana 上看 QPS 曲线、延迟波动。
      • Tracing(链路追踪):全面对接 OpenTelemetry 标准。一条消息从生产者发出、经过 MQ、再到消费者接收,可以在 Jaeger 或 SkyWalking 的分布式链路追踪大盘上被一眼看穿。

通过以上介绍,你会发现无论是 5.x 宏观的 Proxy 存算分离,还是微观的 冷热分层存储,以及全面拥抱 gRPC 协议、强化可观测性标准,RocketMQ 的技术演进终局都是在向 “云原生化” 靠拢——无状态的拼命变轻,有状态的拼命往云端廉价存储上扔。不过,无论数据是存在本地 SSD,还是被偷运到云端 OSS,它在文件里面的组织结构(也就是怎么排列组合的)是绝对一模一样的。


RocketMQ 5.x 的安装和使用

这里的安装和使用,强调的是 “体感”,深入使用的方法,我们会在后面的文章逐步展开。

环境的准备

  • 推荐版本:RocketMQ 5.4.0 稳定二进制包(Binary Release)。
  • 环境依赖:JDK 8 或 JDK 11 以上(5.x 推荐 JDK 11 运行环境)。
  • 物理 facts:RocketMQ 默认 JVM 启动配置是为大厂生产服务器设计的(默认吃 4G~8G 内存)。直接启动,普通电脑或云服务器会瞬间因为内存不足(OOM)而闪退报错。因此,启动前必须修改内存配置!


服务端的简单安装

步骤一:下载与解压

1
2
3
4
5
6
# 1. 下载官方二进制包
wget https://archive.apache.org/dist/rocketmq/5.4.0/rocketmq-all-5.4.0-bin-release.zip

# 2. 解压文件
unzip rocketmq-all-5.4.0-bin-release.zip
cd rocketmq-all-5.4.0-bin-release

步骤二:缩减 NameServer 内存配置

1
2
3
4
5
6
7
8
# vim bin/runserver.sh
# 找到包含 JAVA_OPT 且指定 -Xms 和 -Xmx 的位置(大约在 30-40 行左右),将其改小(本地开发 512M 足够):

# 修改前(示例,不同版本略有差异)
JAVA_OPT="${JAVA_OPT} -server -Xms4g -Xmx4g -Xmn2g ..."

# 修改后(适合本地调测和轻量级生产)
JAVA_OPT="${JAVA_OPT} -server -Xms512m -Xmx512m -Xmn256m"

步骤三:缩减 Broker 内存配置 & 配置文件修改

1
2
3
4
5
6
7
# vim bin/runbroker.sh,修改它的 JVM 内存:

# 修改前
JAVA_OPT="${JAVA_OPT} -server -Xms8g -Xmx8g ..."

# 修改后
JAVA_OPT="${JAVA_OPT} -server -Xms1g -Xmx1g -Xmn512m"

生产核心配置补丁:打开 conf/broker.conf 文件。如果在云服务器或本地多网卡环境下,必须显式指定外网 IP/本地 IP,否则客户端连不上!在文件末尾追加:

1
2
3
4
5
# 显示指定 NameServer 地址,告诉 Broker 去哪台机器的 NameServer 注册。
namesrvAddr=192.168.1.7:9876
# 显式指定 Broker 的宿主机 IP,告诉 NameServer:“我(Broker)的对外服务 IP 是这个,让客户端(Producer/Consumer)通过这个 IP 来连我”
# 在 RocketMQ 的源码和设计中,定义了 brokerIP1 和 brokerIP2,这里的 1 代表 “第一顺位的对外 IP 地址”(或者是主 IP)。
brokerIP1=192.168.1.7

服务启动命令:RocketMQ 的物理时序必须是:先启动通信录(NameServer),再启动大管家(Broker)。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
# 1. 启动 NameServer(常驻后台,并重定向日志)
## 通常不需要重启 mqnamesrv,NameServer 在 RocketMQ 中扮演的是一个无状态的注册中心。它的工作机制是“被动接收”和“定时清理”
## Broker 启动后,会主动向 NameServer 发送心跳包,并把自己最新的 brokerIP1 等配置信息注册上去。
## 即使旧的注册信息还在,NameServer 每隔 120 秒也会自动清理掉连续 30 秒没发心跳的过期 Broker。
$ sh bin/mqshutdown namesrv
$ nohup sh bin/mqnamesrv > ~/logs/rocketmqlogs/namesrv.log 2>&1 &

# 验证 NameServer 是否启动成功(看到 NamesrvStartup 进程)
$ jps

# 2. 启动 Broker(指定刚才修改的配置文件)
$ sh bin/mqshutdown broker
$ nohup sh bin/mqbroker -c conf/broker.conf > ~/logs/rocketmqlogs/broker.log 2>&1 &

# 验证 Broker 是否启动成功(看到 BrokerStartup 进程)
$ jps

常用命令:

1
2
3
4
5
6
7
8
9
# 查看进程
ps -ef |grep rocketmq

# 创建 topic
export NAMESRV_ADDR=192.168.1.7:9876
sh bin/mqadmin updateTopic -n 192.168.1.7:9876 -b 192.168.1.7:10911 -t MyTopic

# 查看 broker 状态
sh bin/mqadmin brokerStatus -n 192.168.1.7:9876 -b 192.168.1.7:10911

端口说明:

  • NameServer:9876,客户端和Broker的注册与发现
  • Broker:10911,消息读写端口
  • Broker:10912,主从同步端口
  • Broker:10909,VIP 通道(可选)
  • Dashboard:8082,web 管理界面


客户端的简单使用

引入依赖

环境起来了,我们来一次无缝的本地代码演练。在你项目中,引入 5.x 兼顾 Remoting 经典协议的客户端 SDK。

1
2
3
4
5
<dependency>
<groupId>org.apache.rocketmq</groupId>
<artifactId>rocketmq-client</artifactId>
<version>5.4.0</version>
</dependency>


生产者代码

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import java.nio.charset.StandardCharsets;

public class OmsProducer {
public static void main(String[] args) throws Exception {
// 1. 初始化生产者,指定业务组名
DefaultMQProducer producer = new DefaultMQProducer("owlias_order_producer_group");
// 2. 绑定刚才启动的 NameServer
producer.setNamesrvAddr("192.168.1.7:9876");
// 3. 启动,和服务端建立连接
producer.start();

System.out.println("====== 生产者启动成功,开始发送订单消息 ======");

for (int i = 1; i <= 10; i++) {
// 4. 构建消息:Topic(大主题), Tag(小标签), Body(业务数据)
Message msg = new Message(
"TOPIC_ORDER",
"Tag_Pay_Success",
("订单数据_ID_00" + i).getBytes(StandardCharsets.UTF_8)
);

// 5. 同步发送
SendResult result = producer.send(msg);
System.out.printf("发送成功 -> 消息ID: %s, 投递队列: %d%n",
result.getMsgId(), result.getMessageQueue().getQueueId());
}

// 6. 释放连接
producer.shutdown();
}
}


消费者代码

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.message.MessageExt;
import java.nio.charset.StandardCharsets;
import java.util.List;

public class OmsConsumer {
public static void main(String[] args) throws Exception {
// 1. 初始化推模式消费者,指定组名(同组分摊消费)
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("owlias_order_consumer_group");
// 2. 指定 NameServer 地址
consumer.setNamesrvAddr("192.168.1.7:9876");
// 3. 订阅 Topic,"*" 代表不过滤,听取该 Topic 下的所有消息
consumer.subscribe("TOPIC_ORDER", "*");

// 4. 注册监听器,Broker 有消息会主动回调这个方法
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
for (MessageExt msg : msgs) {
String body = new String(msg.getBody(), StandardCharsets.UTF_8);
System.out.printf("【消费者接收】监听到新消息! 内容: %s, 来自队列: %d%n",
body, msg.getQueueId());
}
// 5. 核心 Ack:成功消费,告诉 Broker 别再发了
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});

// 5. 启动消费者
consumer.start();
System.out.println("消费者已进入长轮询监听状态...");
}
}


测试验证

先运行 OmsConsumer,你会看到控制台常驻挂起,输出消费者已进入长轮询监听状态。随后运行 OmsProducer 发送三条消息,消费者控制台会立刻瞬间响应并交替打印。

1
2
3
4
5
6
7
8
9
10
11
====== 生产者启动成功,开始发送订单消息 ======
发送成功 -> 消息ID: 24098A002455B220F1DE3A3691F4A4798E3722D8CFE01C186EF30000, 投递队列: 0
发送成功 -> 消息ID: 24098A002455B220F1DE3A3691F4A4798E3722D8CFE01C186F0A0001, 投递队列: 1
发送成功 -> 消息ID: 24098A002455B220F1DE3A3691F4A4798E3722D8CFE01C186F0F0002, 投递队列: 2
发送成功 -> 消息ID: 24098A002455B220F1DE3A3691F4A4798E3722D8CFE01C186F130003, 投递队列: 3
发送成功 -> 消息ID: 24098A002455B220F1DE3A3691F4A4798E3722D8CFE01C186F170004, 投递队列: 0
发送成功 -> 消息ID: 24098A002455B220F1DE3A3691F4A4798E3722D8CFE01C186F1C0005, 投递队列: 1
发送成功 -> 消息ID: 24098A002455B220F1DE3A3691F4A4798E3722D8CFE01C186F220006, 投递队列: 2
发送成功 -> 消息ID: 24098A002455B220F1DE3A3691F4A4798E3722D8CFE01C186F260007, 投递队列: 3
发送成功 -> 消息ID: 24098A002455B220F1DE3A3691F4A4798E3722D8CFE01C186F2A0008, 投递队列: 0
发送成功 -> 消息ID: 24098A002455B220F1DE3A3691F4A4798E3722D8CFE01C186F2D0009, 投递队列: 1
1
2
3
4
5
6
7
8
9
10
11
消费者已进入长轮询监听状态...
【消费者接收】监听到新消息! 内容: 订单数据_ID_007, 来自队列: 2
【消费者接收】监听到新消息! 内容: 订单数据_ID_005, 来自队列: 0
【消费者接收】监听到新消息! 内容: 订单数据_ID_006, 来自队列: 1
【消费者接收】监听到新消息! 内容: 订单数据_ID_009, 来自队列: 0
【消费者接收】监听到新消息! 内容: 订单数据_ID_0010, 来自队列: 1
【消费者接收】监听到新消息! 内容: 订单数据_ID_004, 来自队列: 3
【消费者接收】监听到新消息! 内容: 订单数据_ID_008, 来自队列: 3
【消费者接收】监听到新消息! 内容: 订单数据_ID_002, 来自队列: 1
【消费者接收】监听到新消息! 内容: 订单数据_ID_001, 来自队列: 0
【消费者接收】监听到新消息! 内容: 订单数据_ID_003, 来自队列: 2

注意到后面的来自队列 0/1/2/3 了吗? 这就是我们即将解开的秘密:虽然你的代码感觉是把 10 条消息丢给了同一个 TOPIC_ORDER,但在底层的物理磁盘上,RocketMQ 默认为每个 Topic 创建了 4 个 ConsumeQueue(逻辑消费队列),并在高并发时自动把消息分摊写入。


增加 Dashboard

安装和配置:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
# 下载解压
$ wget https://github.com/apache/rocketmq-dashboard/archive/refs/tags/rocketmq-dashboard-2.1.0.tar.gz
$ tar -zxvf rocketmq-dashboard-2.1.0.tar.gz

# 修改配置文件
$ cd rocketmq-dashboard-rocketmq-dashboard-2.1.0
$ vim src/main/resources/application.yml
# 修改1
rocketmq.config.namesrvAddr=192.168.1.7:9876
# 修改2
## 注释掉 proxyAddr、proxyAddrs(由于我们走的是 Remoting 经典架构,压根没启动 Proxy 进程)

# 编译打包
$ mvn clean package -DskipTests

# 运行
$ nohup java -jar rocketmq-dashboard-2.1.0.jar > /tmp/dashboard.log 2>&1 &

# 查看日志
$ tail -f /tmp/dashboard.log

踩坑点1:源码安装,可能存在跨域问题,修改源码跨域设置:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
// 修改类 org.apache.rocketmq.dashboard.config.AuthWebMVCConfigurerAdapter;

@Override
public void addCorsMappings(CorsRegistry registry) {

registry.addMapping("/**")
// 👉🏻 根据自己的需求设置,修改成自己允许的源
.allowedOriginPatterns("*") // .allowedOriginPatterns("http://localhost:3003")
.allowedMethods("GET", "HEAD", "POST", "PUT", "DELETE", "OPTIONS")
.maxAge(3600)
.allowCredentials(true)
.allowedHeaders("content-type", "Authorization", "X-Requested-With", "Origin", "Accept")
.exposedHeaders("authorization");
}

踩坑点2:如果源码运行报错,可能是 pom 中没有加 lombok

1
2
3
4
5
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<version>1.18.44</version>
</dependency>

擦坑点3:前后端分离 API 设置。编辑 frontend-new/src/api/remoteApi/remoteApi.js

1
2
3
4
5
6
7
8
9
// 修改前
const appConfig = {
apiBaseUrl: 'http://localhost:8082'
};

// 修改后
const appConfig = {
apiBaseUrl: 'http://192.168.1.7:8082'
};


RocketMQ 的存储架构

设计的出发点

在分布式、高并发的消息中间件中,存储面临着两个主要的物理矛盾:

  • 磁盘随机 I/O 是龟速,顺序 I/O 是神速: 传统的机械硬盘甚至普通的固态硬盘(SSD),如果磁头需要不停地在不同的物理文件、不同的磁道之间跳跃擦写(随机 I/O),每秒能抗几百次(IOPS)就到顶了。但如果磁头死守在一个文件的末尾,像磁带一样闭着眼睛一路往后追加写入(顺序 I/O),速度几乎可以媲美内存。
  • 海量 Topic 与文件句柄的冲突:比如 Kafka 的设计是一个 Partition(队列)对应一个物理文件。如果一个电商系统在生产环境建了 5000 个 Topic,每个 Topic 有 4 个队列,磁盘上就会瞬间多出 20000 个物理文件。当高并发大促来临时,十万条消息同时砸过来,操作系统必须在 20000 个文件之间来回切换擦写。顺序写瞬间退化为极度致命的随机写,性能雪崩。

为了撞碎上面这堵墙,RocketMQ 的设计师做出了一个极其大胆且精妙的架构决定:在 Broker 内部,不管你有多少个 Topic,不管消息来自哪里,全部强行按到达的先后顺序,塞进同一个唯一的物理文件里。这个核心文件,就叫做 CommitLog

然而,如果只存这一个文件,消费者在拉取例如上面 TOPIC_ORDER 的消息时,就必须把整个 CommitLog 从头到尾扫描一遍(就像在一本混杂了全校所有学科、所有学生的超级大杂烩笔记本里找某一个人的数学作业),这绝对是不可接受的。


具体涉及的文件

于是,RocketMQ 演进出了大名鼎鼎的三大文件配合矩阵,你可以在服务端 ~/store/ 目录下(默认路径),看到这三个文件夹:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
[root@centos10-01 store]# pwd
/root/store
[root@centos10-01 store]# tree
.
├── abort
#进程死活的“探测仪”,存放 Broker 的进程号。当断电重启,Broker 再次开机时,它第一件事就是去检查有没有 abort 文件。若没有说明上次是正常关机,天下太平。有则说明上次绝对是突发崩溃、强杀或断电。
#Broker 就会立刻启动紧急数据校验(Read-after-Write),顺着 checkpoint 文件去检查 commitlog 末尾有没有写一半的脏字节,并强行将其裁剪修复,死守二进制的安全边界。
├── abort.bak
#这是操作系统或升级时留下的旧探测仪备份。
├── checkpoint
#这个文件极其小,但极其关键。里面记录了三个时间戳:CommitLog 最后一次刷盘时间、ConsumeQueue 最后一次刷盘时间、IndexFile 最后一次刷盘时间。
#当触发崩溃恢复机制时,Broker 就是以 checkpoint 里的时间为基准,去决定到底要从文件的哪一个字节开始重新扫描、修复和分拣。
├── commitlog 真正的消息肉身 👈🏻
│   ├── 00000000000000000000 #这是第 1 个 CommitLog 文件,大小正好是 1GB
│   └── 00000000001073741824 #大小正是 2^30(1G),代表第 2 个文件,文件名就是它在整个 CommitLog 体系里的绝对物理字节偏移量(Offset)
├── compaction
│   └── position-checkpoint
├── config #这个目录下存放的是整个 Broker 的业务运行大脑。它全是 JSON 格式,负责记录所有的行为痕迹
│   ├── consumerFilter.json
│   ├── consumerFilter.json.bak
│   ├── consumerOffset.json # 集群消费的终极肉身!里面密密麻麻地记录着每一个 Consumer Group 在每一个 Queue 文件里分别读到了第几个字节(Offset 点位)。
│   ├── consumerOffset.json.bak # *.json.bak:RocketMQ 每次修改这些元数据前,都会先把旧文件复制一份成 .bak 备份文件,以防写文件时突然断电导致 JSON 损坏,工业级防呆。
│   ├── consumerOrderInfo.json
│   ├── consumerOrderInfo.json.bak
│   ├── delayOffset.json
│   ├── delayOffset.json.bak #记录了老版本 18 个固定级别延迟消息的消费进度。
│   ├── subscriptionGroup.json #记录了当前集群里有哪些消费者组(Consumer Group)以及它们的配置属性。
│   ├── subscriptionGroup.json.bak
│   ├── timercheck
│   ├── timermetrics
│   ├── timermetrics.bak
│   ├── topics.json #记录了当前 Broker 上创建了哪些 Topic,每个 Topic 有几个队列。
│   ├── topics.json.bak
│   ├── transactionMetrics
│   └── transactionMetrics.bak
├── consumequeue 专为消费而生的灵魂索引 👈🏻
│   └── TOPIC_ORDER #这个 TOPIC_ORDER 主题文件夹,点进去正好是 0, 1, 2, 3 四个目录,一个 Queue 就是一个独立的文件夹。
│   ├── 0 #每个文件夹下都有一个 00000000000000000000 文件,里面规规矩矩地躺着每条 20 字节的固定长度索引。
│   │   └── 00000000000000000000
│   ├── 1
│   │   └── 00000000000000000000
│   ├── 2
│   │   └── 00000000000000000000
│   └── 3
│   └── 00000000000000000000
├── index 专为程序员查 dbug 而生的哈希索引 👈🏻
│   └── 20260206104955510 #它的文件名是它的创建时间戳
├── lock
├── rocksdbstore #5.x 的状态元数据引擎,使用嵌入式 KV 数据库 RocksDB。它被用来替代原本简陋的 JSON 文本,专门高速存储 Broker 内部的元数据状态(比如消费队列、轻量级消费位点、或者高级事务状态)。
│   ├── 000012.log
│   ├── CURRENT
│   ├── IDENTITY
│   ├── LOCK
│   ├── MANIFEST-000013
│   ├── OPTIONS-000011
│   └── OPTIONS-000015
└── timerwheel #时间轮——延迟消息的大杀器,在 4.x 版本,延迟消息只支持 18 个固定的延迟级别。而 5.x 版本支持了任意秒数的定时消息!它的底层依托就是这个 timerwheel(时间轮文件)。
  • CommitLog(绝对的肉身,存储真正的消息)
    • 路径:~/store/commitlog/
    • 特点:每个文件默认 1GB。写满了自动创建下一个,文件名以当前文件存储的起始字节偏移量(Offset)命名(例如 00000000000000000000)。
    • 原理:生产者发来的消息(包含 Topic、Tag、Body、各种属性)全量、无脑地顺序追加写入到这里。它是存盘性能快如闪电的根本原因。
  • ConsumeQueue(灵魂索引,专为消费者拉取而生)
    • 路径:~/store/consumequeue/{TopicName}/{QueueId}/
    • 特点:它是轻量级的。每个条目极其固定,只有固定的 20 个字节(8字节的 CommitLog 物理偏移量 + 4字节的消息长度 + 8字节 Tag 的 HashCode)。
    • 原理:当一条消息写入 CommitLog 后,后台会有一个常驻的 “分拣线程”(ReputMessageService),像邮局分拣员一样,一看到 CommitLog 里来新消息了,就立刻把这条消息的“物理地址指针(20字节)”揪出来,顺序投递到对应 Topic 的对应 Queue 文件里。消费者调用 subscribe(“TOPIC_ORDER”) 时,其实读的就是这个极小的索引文件,然后再根据里面的物理地址去 CommitLog 里“指哪打哪”地把真消息瞬间拔出来。
  • IndexFile(哈希索引,专为程序员根据 Key 查消息而生)
    • 路径:~/store/index/
    • 原理:如果在工业上,你想根据业务的 OrderId(或者 RocketMQ 自动生成的 MessageId)去管理后台查某条消息的轨迹。RocketMQ 就会通过底层的 Slot 槽与本地实现的一维 HashMap 结构,把消息的 Key 映射到这个索引文件里。它不影响消费,纯粹是方便排查 Bug 的利器。

结合我们在 Stage 1 写的测试代码,我们可以把上面例子数据流动的大致过程进行梳理:

  1. OmsProducer 连续发送了10条消息(订单 1、2、3 …)。
  2. 这三条消息坐着 TCP 跑车穿过网络,砸进了 Broker。
  3. Broker 的线程拉起这三条消息,顺次排队拍进 commitlog/00000000000000000000 文件的末尾(顺序写盘完成)。
  4. 与此同时,分拣线程开机,发现 CommitLog 多了10 条数据:
    • 第一条是队列 1 的 $\rightarrow$ 在 consumequeue/TOPIC_ORDER/1/ 文件末尾写下 20 字节。
    • 第二条是队列 2 的 $\rightarrow$ 在 consumequeue/TOPIC_ORDER/2/ 文件末尾写下 20 字节。
    • 第三条是队列 3 的 $\rightarrow$ 在 consumequeue/TOPIC_ORDER/3/ 文件末尾写下 20 字节。
    • … …
    • OmsConsumer 拿着长轮询,死守着 TOPIC_ORDER 的各个队列文件。一旦发现里面多了 20 字节的索引,立刻抓出物理偏移量,瞬间去 CommitLog 里把真正的订单 Body 捞出来,触发了你的 consumeMessage 回调函数。这就是为什么我们在上一轮测试中,会看到控制台精准打印出 “来自队列: 0/1/2/3 的底层物理奥秘。”
1
2
3
4
5
6
7
# # 用 hex(十六进制)和 ASCII 码对照的方式查看前 500 个字节
$ hexdump -C /root/store/commitlog/00000000000000000000 | head -n 30

# ConsumeQueue 极其规律,每 20 个字节固定为一条索引
# 20 字节 = 8字节(CommitLog Offset) + 4字节(Size) + 8字节(Tag HashCode)
# 用 od 命令以 20 字节为一组整行对齐打印
$ od -t x8 -N 200 ~/store/consumequeue/TOPIC_ORDER/0/00000000000000000000


RocketMQ 如何做到快写磁盘

在 Java 里,如果我们用传统的 outputStream.write() 或者 FileChannel.write() 向磁盘写一段乱码字节,在 OS 底层会经历一场极其低效的长途跋涉:

  • 用户缓冲区 $\rightarrow$ 内核 PageCache:Java 进程不能直接操作硬件,它得发起一次系统调用(Context Switch,从用户态切换到内核态),CPU 把数据从用户态的 Java 堆内存 拷贝到操作系统内核态的 PageCache。
  • 内核 PageCache $\rightarrow$ 磁盘:由于操作系统不能直接让内存去碰硬盘,通常还要经过 DMA(直接内存存取)把数据从 PageCache 拷到块设备驱动的缓冲区,最后再硬刻进磁盘。
  • 在这个过程中,数据被来回复制了 4 次,上下文切换了 4 次。在高并发、每秒几十万条消息砸过来时,CPU 光用来做搬运工和状态切换就彻底累瘫了。

RocketMQ 存储的大杀器,就是 Java 中的 MappedByteBuffer(底层对应的就是操作系统的 mmap 系统调用)。它的核心思想是:利用操作系统的虚拟内存地址空间,把磁盘上的一个物理文件(比如 1GB 的 CommitLog),直接映射到进程的用户内存地址空间里。通俗来说,操作系统在内存里画了一个圈,对 Java 进程说:“从现在开始,你把这块内存当成那个 1GB 的 CommitLog 文件就行了,你尽管闭着眼睛往这块内存里写乱码字节,剩下的交给我。”

当你的 OmsProducer 发来一条消息,Broker 往 CommitLog 里写盘时,Java 代码实际上是在直接操作这块被映射的内存(Direct Memory)。数据直接写进了操作系统的 PageCache,中间不需要经过用户态到内核态的 CPU 搬运。对于 Java 进程而言,它只要把字节拍进内存,这次写盘操作就“结束”了!整个过程只有 1 次 异步的 DMA 硬件拷贝,CPU 得到了彻底解放。

mmap 只是给了我们像操作内存一样方便的权力,而 RocketMQ 能把这个权力发挥到极致,是因为它严格遵守了顺序写。CommitLog 文件大小严格固定为 1GB。当消息源源不断过来时,RocketMQ 在物理内存里维护了一个自增的物理写入点位(WrotePosition)。每次写消息,就是一条紧挨着上一条,在内存里疯狂做数组追加(memcpy)。操作系统的物理磁盘磁头甚至不需要转动寻道,只需要像磁带机一样一往无前地顺序擦写。 这种顺序写的速度,在现代物理架构下,几乎和纯内存操作没有本质区别。

你可能会有疑问,数据只是写到了 OS 的 PageCache 内存里,Java 进程就认为成功了。如果这时候服务器突然被人拔了电源,这部分还没刻进磁盘磁道的消息不就彻底凭空消失了?这就是 RocketMQ 在高性能与高可靠之间,必须做出的 “刷盘策略抉择”

  • 如果要异步刷盘:消息写入 PageCache 后立刻返回成功。后台常驻线程每隔 500ms 定时调用一次 fsync(),把内存数据真正冲刷到硬盘上。可以做到吞吐量极高,基本可以无视磁盘物理限制。但是有风险(断电瞬间可能丢失过去 0.5 秒的数据,但进程挂掉/kill -9 不丢数据,因为 OS 还在)。RocketMQ 之所以能抗住双十一百万洪峰,秘诀就在于它用 mmap 把“写磁盘”的动作降维打击成了“纯写 PageCache 内存”,再利用异步刷盘策略让磁盘在后台慢慢干活。
  • 如果是同步刷盘:消息写入 PageCache 后,线程必须原地死等 fsync() 强制让操作系统把这批字节刻进磁盘磁道后,才给客户端返回成功。那么吞吐量必然下降,但是安全性极高,真正做到了物理不丢,哪怕拔电源也安然无恙。通常用于金融扣款、核心订单、关键对账的场景。


RocketMQ 延迟消息的本质

在很多普通中间件里,高级特性(如延迟、事务)往往要在架构上加极其复杂的外部插件或分布式锁。而 RocketMQ 神奇的地方就在于:它纯粹依靠在底层三大文件(CommitLog、ConsumeQueue)上玩弄“障眼法”和“改写指针”,就以极度轻量、高性能的姿势玩转了这些高级特性。我们首先来看使用最频繁、设计最绝妙的第一个特性——延迟消息(Scheduled Message)。

如果让你设计延迟消息,最直觉的想法可能是:消费者要 10 分钟后才看得到这条消息,那我就把消息在内存里放 10 分钟,或者在磁盘里建个定时器,10 分钟到了再写入 CommitLog,但如果突发 100 万条延迟消息,内存直接爆掉;如果用定时器频繁扫描普通磁盘文件,顺序写会瞬间退化为致命的随机读写,性能雪崩。RocketMQ 的破局智慧叫做:偷梁换柱。

生产者发送:隐形斗篷

当你的 Java 代码发送了一条延迟级别为 3(例如 10 秒延迟)的订单消息,目标主题是 TOPIC_ORDER。

  • 第一步:无脑落盘。Broker 收到消息后,不管三七二十一,先把消息原本、真实地顺序写入 CommitLog。
  • 第二步:暗度陈仓(核心黑魔法)。在消息即将投递到 ConsumeQueue 的前一刹那,分拣线程ReputMessageService 对这条消息动了外科手术:
    • 它把消息的真实主题 TOPIC_ORDER 和队列 ID 偷偷揪出来,塞进消息的属性(Properties)里藏好。
    • 它把这条消息的真实主题改写为了一个系统内置的秘密主题:SCHEDULE_TOPIC_XXXX。
    • 它根据你的延迟级别(级别 3),把消息投递到了对应的第 2 号队列(QueueId = 延迟级别 - 1)。

因为我们的 OmsConsumer 监听的是 TOPIC_ORDER,而此刻这条消息的索引躺在 SCHEDULE_TOPIC_XXXX 文件里,消费者在根本看不到它。这就完美实现了 “延迟挂起” 的效果,且依然保持了 CommitLog 的绝对顺序写!


守护者:时间到现原形

在 Broker 内部,针对系统自带的 SCHEDULE_TOPIC_XXXX 的每一个延迟级别队列,都安排了一个专门的常驻守护线程(每个队列一个)。

  • 这个线程像一个不知疲倦的钟表,死死盯着对应延迟级别的 ConsumeQueue 文件。
  • 它每次只读下一个 20 字节的索引,根据索引去 CommitLog 把消息捞出,看看里面的“投递时间”到了没有。
  • 如果时间没到,线程原地睡眠(比如睡 100ms)然后继续看;
  • 如果时间到了(10 秒过去):线程会把这条消息重新拽出来,把藏在属性里的真实主题 TOPIC_ORDER 重新恢复到消息头上,然后再次、原封不动地重新写入一遍 CommitLog!
  • 这一次,由于主题已经是 TOPIC_ORDER 了,正常的分拣线程看到后,就会把索引规规矩矩地投递到 consumequeue/TOPIC_ORDER/ 下。消费者长轮询瞬间命中,回调你的 Java 业务代码!


这样设计的优点

通过这种 “先改名藏起来、到期了再改名写回来” 的移花接木术,RocketMQ 达到了极其恐怖的工业性能:

  • 零内存开销:所有的延迟消息都在磁盘文件里顺序躺着,哪怕堆积一个亿,也绝不吃系统内存。
  • 绝对的顺序性:无论是藏进去,还是到期拿出来,由于全部统一在 CommitLog 尾部追加,磁盘磁头不需要做任何无意义的来回摆动(随机 I/O)。


RocketMQ 事务消息的实现

在传统的分布式事务(如 XA 模式、二阶段提交)中,系统会通过高成本的 “强锁” 机制让所有参与者进入等待状态。这种做法在互联网高并发下无异于自杀,任何一个节点卡顿都会拖垮整条链路。而 RocketMQ 的事务消息玩了一场更加高明的游戏:它在底层文件系统里,依然沿用了延迟消息类似的 “移花接木” 与 “半消息(Half Message)” 障眼法,在完全不锁数据库、不锁网络通道的前提下,优雅地实现了可靠消息的最终一致性。为了让你彻底看清它在磁盘里的肉身,我们把整个发消息、扣款、提交的宏观过程,拆解为底层的三大文件流动时序。


第1阶段:发半消息

当客户端调用 sendMessageInTransaction 发出一条消息(比如:订单已支付,通知积分系统加 100 积分):

  • Broker 收到消息,老规矩,先把消息原本、真实地顺序写入 CommitLog。
  • 在分拣线程准备写 ConsumeQueue 索引时,发现这是一条事务半消息。于是,它再次施展“换头术”:
    • 把原本的目标主题(如 TOPIC_POINTS)偷偷藏进属性。
    • 把消息的主题强制改写为系统内置的:RMQ_SYS_TRANS_HALF_TOPIC。
    • 将索引顺序投递到这个内置主题的 ConsumeQueue 里。
    • 由于积分系统的消费者只监听了 TOPIC_POINTS,这条消息在 HALF_TOPIC 里处于完全隐形的状态,消费者根本看不到,这就是 “半消息”。


第2阶段:执行本地事务

此时 Broker 腾出手来,给你的 Producer 异步返回一个 Ack。Producer 收到后,立刻在本地启动你的业务代码:执行本地数据库扣款/订单更新。

  • 情况 A(皆大欢喜):本地数据库成功执行,Producer 给 Broker 回复一个 COMMIT 信号。
  • 情况 B(本地崩溃):本地数据库执行失败(比如余额不足抛出异常),Producer 给 Broker 回复一个 ROLLBACK 信号。


第3阶段:事务最终裁决

当 Broker 收到 Producer 发来的最终裁决(Commit 或 Rollback)时,它是怎么操作磁盘文件的?注意:因为 CommitLog 只能顺序追加,RocketMQ 绝对不会、也无法去修改或删除第一阶段写进去的那条半消息!

  • 如果收到的是 COMMIT:
    • 从物理上来说,那条半消息依然永远躺在 CommitLog 过去的某个角落里。
    • Broker 会把这条半消息的躯壳重新拽出来,把属性里藏着的真实主题 TOPIC_POINTS 恢复到消息头上。
    • 将这条消息重新、完整地顺序追加写入一遍 CommitLog 的末尾!
    • 此时,正常的分拣线程看到它是一条正常消息,规规矩矩地在 consumequeue/TOPIC_POINTS/ 下写下 20 字节索引。积分消费者长轮询瞬间命中,开始加积分。
  • 如果收到的是 ROLLBACK:
    • Broker 表现得极其冷漠。它什么都不做,不重写消息,也不生成目标 Topic 的索引。
    • 随着时间推移,那条躺在 HALF_TOPIC 里的半消息,会被操作系统的定时清理机制连同旧文件一起无情抹去。积分系统永远不会知道这条消息存在过,两端逻辑完美保持干净。


突然断电了怎么办?

如果 Producer 执行本地事务时,突然应用崩溃、断电、或者网络断开,导致 Broker 迟迟没有收到 COMMIT 或 ROLLBACK,这条 “隐形” 的半消息难道要在磁盘里躺一辈子吗?RocketMQ 在底层设计了一个常驻的事务回查线程(TransactionCheckService):

  • 它在后台默默扫描系统内置的 RMQ_SYS_TRANS_HALF_TOPIC 队列文件。
  • 如果发现某条半消息已经在里面躺了超过 15 秒(可配置),且既没有被 Commit 也没有被 Rollback。
  • Broker 会主动发起一次物理反击(回查):它顺着网络通道,向你的 Producer 集群随机挑选一台机器发送一个“回查请求”:“哥们,15 秒前那笔订单,你们本地数据库到底扣款成功了没有?”
  • 你的 Java 业务代码(实现 executeLocalTransaction 接口)会去查一下本地数据库是否有这条记录,并把结果再次回复给 Broker。
  • 拿到结果后,Broker 再次重复上面的真假提交逻辑。回查上限默认 15 次,如果 15 次都找不到人,这条消息才会被彻底丢进死信队列(DLQ)。

RocketMQ 事务消息的本质:通过向一个消费者看不见的内置“停机坪”(HALF_TOPIC)顺序追加写,等本地事务尘埃落定后,再通过在 CommitLog 末尾重新追加写(Commit)或者直接装死(Rollback)的手段,用极其轻量级的“顺序写盘”完美平替了重型的分布式锁协议。


RocketMQ 如何做到消息不乱

在分布式高并发环境下,要保证消息的绝对顺序是一个极其反直觉的物理难题。如果你的业务有三步:创建订单 $\rightarrow$ 支付订单 $\rightarrow$ 发货订单。当百万 QPS 洪峰砸过来,多线程并发发送,网络抖动,如果不对底层做物理限制,发货消息完全可能比创建消息先到达 Broker,或者消费者端的多个线程同时拉到了这三条消息,导致“先发货、再创建”的业务灾难。

RocketMQ 解决这个问题的铁律极其纯粹:局部有序(分区有序)。它在底层三大文件上不玩任何花哨的障眼法,而是通过硬核的 “通道绑定” 与 “物理单线程化”,实现了顺序闭环。


发送端:单行道发送

在 4.x/5.x 经典架构下,当你创建一个 Topic(如 TOPIC_ORDER)时,Broker 默认会为它在磁盘上创建 4 个 Queue(队列 0, 1, 2, 3)。这意味着在本地 ~/store/consumequeue/TOPIC_ORDER/ 下,天然存在 4 个独立的文件夹。如果采取普通发送模式,RocketMQ 会使用轮询算法(RoundRobin),把你的创建、支付、发货消息分别扔进 Queue 0、Queue 1、Queue 2,那么消费者在多线程并发拉取时,由于不同队列的网络和处理速度不同,顺序必然乱套。

为了保证顺序,你的 Java 业务代码在发送时,必须使用哈希选择器,强行把具有相同 OrderId 的所有消息,砸进同一个固定的 Queue 文件里。

1
2
3
4
5
6
7
8
9
10
// 确保同一个订单号的消息,永远走同一个队列文件
// ::局部有序的物理铁律一:在同一个 ConsumeQueue 文件里,消息索引在物理上永远是绝对有序的。
SendResult sendResult = producer.send(msg, new MessageQueueSelector() {
@Override
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
Long orderId = (Long) arg;
long index = orderId % mqs.size(); // 根据订单号取模
return mqs.get((int) index); // 强行返回固定的队列
}
}, orderId);

通过这一步,同一个订单的 “创建 $\rightarrow$ 支付 $\rightarrow$ 发货” 三条消息,在网络传输的尽头,会规规矩矩、按先后顺序排队拍进同一个 CommitLog 的末尾。随后,分拣线程会将它们的 20 字节索引,完全顺序地、一条挨着一条写进同一个队列文件(比如 consumequeue/TOPIC_ORDER/queue_1)中。

需要注意的是,单行道消息发送的时候一定是使用的是 MessageQueueSelector。如果是普通的消息发送,根本不会有单行道发送的效果:

1
2
3
4
5
6
7
8
9
// 当你直接调用 producer.send(msg) 时,使用的是默认的 TopicPublishInfo 路由逻辑。
// 在这个过程中,底层负责网络分发的线程根本不看你的 Key 长什么样。
// 你的 "ORDER_KEY_10086" 只是作为消息主体属性的一部分,跟着二进制字节流一起顺序写入了对应 Queue
// 所在 Broker 的 CommitLog 文件。随后被分拣线程提取出来,存进了 index/(索引文件) 里。
// 默认模式下,Key 的作用是用来“查”的(基于 IndexFile 高速检索),而不是用来“路由分区”的。
for (int i = 0; i < 100; i++) {
Message msg = new Message("TOPIC_ORDER", "TagA", "ORDER_KEY_CONST", ("集群测试数据:" + i).getBytes());
SendResult sendResult = producer.send(msg);
}

改为单行道发送:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
for (int i = 0; i < 100; i++) {
String orderId = "10086"; // 你的业务 Key
Message msg = new Message("TOPIC_ORDER", "TagA", "ORDER_KEY_" + orderId, ("集群测试数据:" + i).getBytes());

// 核心改造:显式传入 Selector,并将 orderId 作为参数传进去
SendResult sendResult = producer.send(msg, new MessageQueueSelector() {
@Override
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
String id = (String) arg;
// 1. 利用哈希算法,让相同的 id 算出绝对相同的绝对哈希值
int hashCode = id.hashCode();
// 2. 对当前集群所有可用的队列总数(这里是 8)取模
int index = Math.abs(hashCode) % mqs.size();
// 3. 返回固定的队列
return mqs.get(index);
}
}, orderId);

System.out.printf("消息投递成功,此时对应的物理队列是: %s%n", sendResult.getMessageQueue());
}


消费端:两道互斥锁

消息在磁盘文件里排好队了,但消费端通常是多线程并发干活的。如果消费端的线程 A 抓取了 “创建订单”,线程 B 抓取了“支付订单”,线程 B 跑得比 A 快,业务依然会雪崩。为了在消费者端也维持绝对的单行道,RocketMQ 祭出了两道重型物理锁:

第一道锁:分布式锁(Broker 端锁 Queue)

当你的消费者集群启动时,会有多台机器同时去分摊消费 4 个队列。RocketMQ 的底层策略是,任何一个 ConsumeQueue 文件,在同一时刻,只能被集群里的一台消费者机器的某一个线程独占绑定。消费者会定期向 Broker 发送锁请求(lockBatchMQ)。Broker 会在内存里给这个 Queue 加上一把分布式锁。只有抢到这把锁的消费者,才有资格去拉取该 ConsumeQueue 里的二进制字节。这保证了多台机器之间不会产生交叉抢夺。

第二道锁:本地内存锁(线程锁 ProcessQueue)

当消费机器把消息拉回本地内存后,准备分发给本地的线程池(ThreadPoolExecutor)去回调你的 consumeMessage。如果用的是普通监听器(MessageListenerConcurrently),线程池里的 20 个线程会一拥而上瓜分这批消息,顺序再次完蛋。所以顺序消费,必须改用 MessageListenerOrderly(顺序监听器):

1
2
3
4
5
6
7
8
9
// 顺序消费核心骨架
consumer.registerMessageListener(new MessageListenerOrderly() {
@Override
public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs, ConsumeOrderlyContext context) {
// 本地单线程或加锁消费逻辑
System.out.println(Thread.currentThread().getName() + " 顺序消费: " + new String(msgs.get(0).getBody()));
return ConsumeOrderlyStatus.SUCCESS;
}
});

在 MessageListenerOrderly 内部,针对每一个拉取到本地的队列快照(ProcessQueue),都内置了一把 Java 的 ReentrantLock 重入锁。当本地线程池企图分发这批消息时,线程必须先成功抢到这把本地内存锁,才能消费该队列里的消息。由于锁的排他性,这就强制逼迫本地线程池:对于同一个 ConsumeQueue 文件出来的消息,本地只能有一个线程在那里像吃面条一样,一根一根、按磁盘里的物理顺序一口口吃下去。

我们用一句话把 “创建 $\rightarrow$ 支付 $\rightarrow$ 发货” 的路径闭环:

  • Producer 用 Hash 把同一个订单的消息强行砸进同一个 ConsumeQueue 文件(物理排队);
  • Broker 用分布式锁确保这个文件同一时刻只能被一台机器拉取(单机独占);
  • Consumer 用本地内存锁确保本地只有一个线程能顺序消费该文件(单线程执行)。


RocketMQ 如何做到消息不漏

消息不漏的关键设计

消息丢失可能发生在三个阶段:发送时、Broker 存储时、消费时。RocketMQ 在每个阶段都设了物理关卡。

  • 发送阶段:同步/异步的 ACK 审判

    • 当 OmsProducer 发送消息后,底层 Netty 会原地死等(或者通过回调机制等待)Broker 返回的 SendResult。
    • 如果网络丢包或 Broker 挂了,Producer 会在超时后自动触发重试机制(默认重试 2 次,且会自动切换到其他健康的 Broker 队列)。
    • 只有当收到包含 SendStatus.SEND_OK 的响应报文时,发送端才宣告“此消息未丢”。
  • 存储阶段:同步刷盘 + 物理多副本

    • 消息到了 Broker 的 PageCache 内存后,同步刷盘(物理不丢),强制调用 fsync() 将字节彻底刻进磁盘磁道后才返回 SEND_OK。
    • DLedger/主从复制(架构不丢):在企业级部署中,Master 收到消息后,会通过 Netty 将其同步或异步复制到 Slave 机器。即使 Master 所在的物理机瞬间爆炸,Slave 也能无缝接管,数据在多台机器的磁盘上依然完好无损。
  • 消费阶段:提交位点(Offset)的受控右移

    • 当消费者把消息拉到本地后,如果在执行你的 consumeMessage 业务代码时突然断电,那么 RocketMQ 绝对不会在 “拉取消息” 时就认为消费成功。只有当你的业务代码规规矩矩地返回 ConsumeConcurrentlyStatus.CONSUME_SUCCESS 时,消费者才会向 Broker 发送 UPDATE_CONSUMER_OFFSET 请求。
    • 只要没有返回成功,Broker 端的 ConsumeQueue 消费指针(Offset)就绝对不会向右移动。机器重启后,会顺着上一次没提交的点位重新拉取,确保不漏。


消费失败后的自愈链路

另外,这里我想讨论一下消费失败后 RocketMQ 的自愈链路,这是一个包含 3 个特殊 Topic 变换的绝妙设计。

1
ACK ⇄ 重试队列(Retry) ⇄ 死信队列(DLQ)
  • 关卡一:拒绝提交与 ACK
    • 当你的消费监听器最终返回了 ConsumeConcurrentlyStatus.RECONSUME_LATER(稍后重试),或者直接抛出了 Throwable 异常。消费端会拒绝向 Broker 提交当前点位。
    • 随即消费端通过 Netty 向 Broker 发送一条 CONSUMER_SEND_MSG_BACK 的特殊报文,把这条惹祸的消息原封不动地退还给 Broker。
  • 关卡二:重试队列(%RETRY%)的梯度阻尼
    • Broker 收到退回消息后,绝对不会原地死循环投递(那会瞬间卡死消费线程)。它再次玩起了“换头术”
    • 它把消息的主题改写为:%RETRY% + ConsumerGroup_Name(重试主题)。
    • 这个重试主题在磁盘里的肉身,本质上就是一个延迟消息 Topic。
    • 消息会根据重试次数(第 1 次重试对应延迟级别 3,即 10 秒后;第 2 次 30 秒后…),被拍进对应的延迟队列里。
    • 10 秒钟过去,延迟到期,后台线程把消息剥离出来,再次送回 %RETRY% 队列中。消费者组其实暗中监听了这个 %RETRY% 主题,于是重新拉取,触发第二次业务尝试。这就是高可靠的梯度退避重试机制。
  • 关卡三:死信队列(%DLQ%)的终极火化
    • 如果老天不作美,你的业务系统整整崩溃了一天。这条消息在重试队列里惊人地重试了 16 次(默认上限),耗时近 4-5 个小时。
    • 当第 17 次砸过来依然失败时,RocketMQ 认为这条消息已经无药可救(可能格式天生写错了)。
    • Broker 会在 CommitLog 末尾最后一次重写该消息,将其主题改写为:%DLQ% + ConsumerGroup_Name(死信主题,Dead Letter Queue)。
    • 一旦进入 %DLQ%,消息在底层的路由将彻底死掉。原有的消费者再也不会去拉取它。它在磁盘里静静地躺着,等待运维人员在 Dashboard(控制台)上手动点击 “Resend Message”(重发死信),或者等待磁盘过期被物理擦除。


RocketMQ 如何做到消息不重

分布式架构在理论上是无法做到绝对的 “不重复” 的!无论是 RocketMQ、Kafka 还是 RabbitMQ,在遭遇网络极限抖动时,都必须在“不漏”和“不重”之间做二选一。而消息中间件的原则永远是:宁可重复,绝不能丢失(At least once)。为什么一定会产生重复呢?原因是:

  • 发送端网络抖动:Producer 把消息成功写入了 CommitLog,Broker 刚准备给 Producer 发 ACK 确认,结果局域网断网了。Producer 迟迟拿不到 ACK,认为发送失败,于是重新发送了一遍。这时候 CommitLog 里就躺了两条一模一样的消息。
  • 消费端网络抖动:消费者把消息完美执行,业务数据库也落盘了,正准备给 Broker 提交消费点位(Offset)时,消费者机器被 kill -9 强杀了。Broker 认为这批消息没消费成功,换了一台消费者机器重新投递。

既然消息中间件在物理上无法避免重复投递,那么“不重”的最后一公里,必须在开发者的 Java 业务代码中完成。这就是大名鼎鼎的 消费幂等。我们在生产上最常用的三大物理防御手段:

  • 数据库唯一索引(Unique Key):如果消费的是订单创建,利用 order_id 作为数据库唯一索引。第二次重复消息砸过来,数据库会直接抛出 DuplicateKeyException,业务直接捕获并优雅返回成功即可。
  • Redis 分布式布隆过滤器/去重表:消息一过来,先拿着消息的唯一标识(如 MessageId 或业务 BizKey)去 Redis 里 setnx。如果返回失败,说明之前已经吃过这条消息了,直接丢弃。
  • 状态机幂等:针对 “支付 $\rightarrow$ 发货” 业务,数据库记录里订单状态已经是 “已发货” 了。如果又来了一条 “支付” 的消息,代码一看状态已经往后演进了,直接宣告忽略。


集群消费和广播消费

在 Java 代码中,我们通过 setMessageModel 来切换这两种消费模式。它们在底层的差异可以用一句话概括:消费点位(Offset)的肉身到底存在哪。

1
2
3
consumer.setMessageModel(MessageModel.CLUSTERING); // 集群模式(默认)
// 或
consumer.setMessageModel(MessageModel.BROADCASTING); // 广播模式


集群消费(Cluster)

总结起来就是 “一干多帮,点位上云”。比如:

  • 一个 Consumer Group(例如 GID_ORDER)下面有 3 台机器。一条订单消息过来,只能被其中一台机器吃掉。
  • 因为大家合伙干活,消费进度必须共享。因此,它们的消费进度(Offset)是集中存储在 Remote Broker 端的(对应物理文件 ~/store/config/consumerOffset.json)。
  • 每次某台机器消费完,都会通过 Netty 向 Broker 汇报:“queue_1 我已经读到第 500 字节了”。Broker 统一记录,其他机器扩容进来时,直接去 Broker 问点位,接着干活。


广播消费(Broadcast)

总结起来就是 “各过各的,点位下沉”。

  • 每台机器都必须收到全量的消息(通常用于本地内存缓存刷新、网关推送)。
  • 由于每台机器的消费节奏完全独立,Broker 根本无法也无需替它们管理点位。它们的 Offset 死死地保存在消费者本地的磁盘上(位于消费者机器的 ~/.rocketmq_offsets/... 纯文本文件里)。
  • 如果某台机器被 kill -9 重启且本地文件被删,它将彻底丢失自己的进度,只能从头或者从最新的地方开始裸奔。


负载均衡与数量匹配

在集群模式下,一个 Topic 的 4 个 ConsumeQueue 文件夹,是如何在 3 台消费者机器之间分赃的?这就是 AllocateMessageQueueStrategy(队列分配策略)的职责。当消费者启动或扩容时,每台机器都会在本地执行一套完全相同且无锁的数学矩阵算法(默认是 AllocateMessageQueueAveragely 平均分配算法)。无论哪种负载均衡,必须要遵循队列数和消费者数匹配的原则。

假设你的 Topic 默认有 4 个 Queue,我们在集群里部署不同数量的 Consumer 机器,看看底层的物理分配事实:

  • 事实 1:Queue(4) ⇄ Consumer(4) —— 完美平衡(1:1) —— 每台机器分到 1 个 Queue 文件夹的读权限。磁头各司其职,网卡流量完美均摊,这是最健康的工业状态。
  • 事实 2:Queue(4) ⇄ Consumer(2) —— 减半承载(2:1) —— 每台机器分到 2 个 Queue。每台机器本地会启动 2 个长轮询拉取线程分别死守这两个文件,虽然压力翻倍,但流量依然是均衡的。
  • 事实 3:Queue(4) ⇄ Consumer(5) —— 严重的致命浪费(4:5)—— 当 Consumer 数量大于 Queue 数量时,神奇的无锁算法会导致:
    • Consumer 1 $\rightarrow$ 分到 Queue 0
    • Consumer 2 $\rightarrow$ 分到 Queue 1
    • Consumer 3 $\rightarrow$ 分到 Queue 2
    • Consumer 4 $\rightarrow$ 分到 Queue 3
    • Consumer 5 $\rightarrow$ 分不到任何 Queue!只能原地躺平,变成 “僵尸节点”。

所以,这就告诉我们:

  • 一个 ConsumeQueue 文件在同一时刻只能被同组内的一个 Consumer 线程独占。
  • 工业上如果发现消费能力不足想要扩容机器,第一步不是盲目加 Consumer 机器,而是必须先去控制台把 Topic 的 Write/Read Queue 数量改大,否则加进去的机器也是白白烧电!

RocketMQ 网络通讯架构

如果说 CommitLog、ConsumeQueue 是 RocketMQ 的躯干和肢体,那么基于 Netty 构建的高性能网络传输集群就是它敏捷强悍的 “神经系统”。在百万 QPS 的极端高并发下,普通的 I/O 模型早就因线程阻塞、上下文切换或内存拷贝而吐血身亡。RocketMQ 能做到气定神闲,完全得益于它对 Netty 的深度定制,以及一套精妙绝伦的 Reactor 多线程接力模型。

Broker 的三层 Reactor

当一条消息的字节流顺着 TCP 网线砸进 Broker 主机(比如 192.168.1.7)时,在 Broker 内部,它并不是直接由一个线程一包到底去写盘的,而是像工厂流水线一样,经历了一场极其严密的线程接力赛。RocketMQ 在底层设计了三层 Reactor 线程,把 “连接”、“读写”、“分发” 这三件事拆得清清楚楚:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
[ 网线字节流 ] 


1. EventLoopGroupBoss (1个线程) ─────── 专职:接电话(Accept 建立 TCP 连接)


2. EventLoopGroupSelector (默认32个) ── 专职:收发信(网卡字节流 ⇄ Netty 报文结构)


3. NettyServerCodec (SSL/编解码) ───── 专职:翻译官(二进制字节 ⇄ RemotingCommand 对象)


[ NettyServerWorkerThread ] ────────── 专职:保安分拣(把请求精准扔进对应的业务线程池)

  • BossGroup(1个线程):死守端口(如 10911)。别的啥也不干,专门负责监听客户端的连接请求。一旦建立连接,立刻把这条长连接 SocketChannel 丢给 SelectorGroup,自己转头继续等下一个连接。
  • SelectorGroup(网络 I/O 线程池,默认 32 个线程):基于 Linux 的 epoll 机制,专门负责在这堆长连接管道里轮询读取网络字节流、或者把响应字节流推回给客户端。
  • 编解码处理器(NettyServerCodec):当 SelectorGroup 抓到一堆二进制乱码后,传给这个处理器。它根据 RocketMQ 独有的协议头,把字节流翻译成程序员能看懂的 RemotingCommand(通讯命令对象)。


八大业务线程池

翻译成 RemotingCommand 后,按理说该轮到写磁盘了。但是 Selector 线程是绝对不能用来写磁盘的(因为 mmap 再快,也会有短暂的 PageCache 锁竞争,一旦 I/O 线程被阻塞,整个 Broker 的网络吞吐瞬间归零)。于是,RocketMQ 再次解耦,在 Netty 之后安排了 8 大业务线程池。Netty 线程会根据请求的 “功能类型”(RequestCode),把对象像接力棒一样塞进对应的线程池里。


消息网络通讯全景图

我们结合之前的示例,以生产者写消息为例,感受一下一条消息在 Netty 线程模型里的终极丝滑旅行:

  • 你的 Java 线程调用 send(),Netty 客户端把消息打包通过 TCP 砸向服务器。
  • Broker 的 Boss 线程接收连接,Selector 线程顺着管道疯狂吸入字节流。
  • Codec 处理器手起刀落,把字节流翻译成一个 RequestCode 为 SEND_MESSAGE 的 RemotingCommand。
  • Netty 线程看了一眼类型,反手把对象精准抛进了 SendMessageExecutor 线程池,Netty 线程立刻抽身去接下一个网络请求。
  • SendMessageExecutor 里的一个线程接过接力棒,拿着这个对象,顺着虚拟内存地址,一巴掌把消息追加写入 mmap 映射的 CommitLog 内存(PageCache)里。
  • 写完后,线程立刻拼装一个 SEND_OK 的响应对象,再次丢给 Netty 的 Selector 线程,顺着网线飞回你的 OmsProducer。

这就是 RocketMQ 承载百万吞吐的网络通讯底牌:网络的归网络(Netty),业务的归业务(8大线程池),层层递进,绝不阻塞。


RocketMQ 的时间轮算法

最后,我们来看一个 RocketMQ 中比较轻松有趣的小工具——时间轮。在 store/ 目录下,任意秒定时消息之所以能跑起来,靠的是两个文件的物理配合:

  • timerwheel(时间轮文件):充当钟表的刻度盘。就是一个严格等分的物理圆环。这个表盘被切分成了很多个小格子(Slot槽),每一个格子代表现实世界里的 1 秒钟。比如 RocketMQ 的时间轮是严格的 2 天的秒数(也就是 $2 \times 24 \times 3600 = 172800$ 个格子)。每个格子在磁盘里的空间是严格固定(比如 32 字节)的。因为大小固定,时间轮文件的大小也是固定的。指针每走一秒,就对应磁盘文件里往后移动 32 字节的物理位置。走到末尾,又绕回文件开头(取模循环)。
  • timerwheel/../timerlog(定时日志文件):充当 “记事本/挂历”。因为每秒钟可能有很多条定时消息到期(比如双十一零点有 10 万条消息同时到期),但时间轮的一个格子(Slot)只有 32 字节,根本塞不下这么多消息。于是,RocketMQ 准备了第二个文件叫 timerlog。只要来了一条定时消息,Broker 就把它顺序追加写进 timerlog 的末尾(跟 CommitLog 一样高效)。timerlog 的每条记录都会记录:消息在 CommitLog 的位置、到期时间、以及指向同一个格子前一条消息的物理指针。

我们用一个具体的例子来看看它们是怎么物理配合的。假设现在是 10:00:00,指针正指在时间轮的 A格子 上。此时,用户发来了一条消息:“请在 10:00:02(2秒后)帮我扣款”。

  • Broker 计算发现,10:00:02 对应时间轮上的 C格子。
  • Broker 先把这条消息写进 timerlog 文件的末尾(假设写在物理位置 Offset = 800 的地方)。
  • 接着,Broker 去看一眼 C格子 里面当前存的是什么。如果 C格子 目前是空的,Broker 就把 C格子 里的数值改成 800。
  • 如果过了 0.1 秒,又来了一条 10:00:02 到期的消息。Broker 把它写进 timerlog 的 Offset = 900 处。然后,把这条新记录的 “前驱指针” 指向 800,再把 C格子 里的数值更新为 900。
  • 通过这种操作,时间轮的每个格子,其实都成了一条隐形链表的表头。格子只保存最新一条消息在 timerlog 里的地址,剩下的消息在 timerlog 内部通过指针像烤串一样串在一起。

接下来就是让这个时间轮动起来了。Broker 内部有一个专门的网卡式长轮询线程,每隔一秒钟,就把时间轮的指针往前拨一个格子。

  • 时间走到了 10:00:02,指针指向了 C格子。
  • 线程一审问,发现 C格子 里的值是 900(不为空)。
  • 线程立刻拎着 900 这个地址,冲进 timerlog 文件里把第 900 字节的消息捞出来:“时间到了,现出原形!”,然后把它重新写回 CommitLog 投递给消费者。
  • 顺着第 900 字节里记录的指针,线程又摸到了 Offset = 800 的前一条消息,拉出来继续投递。
  • 直到拉出的指针为零,说明 C格子 这一秒的所有消息全部处理完毕。

普通的内存时间轮(比如 Netty 内部的)最大的痛点是吃内存。如果用户定了一个 1 天后发送的大消息,消息积压在内存里,一断电全没了。RocketMQ 的磁盘时间轮之所以厉害,就在于它把数据全下沉到了磁盘:

  • 不吃内存:指针移动和链表回溯,全是在 mmap 映射的磁盘文件里倒腾,内存里只留极少的高频热点缓存。
  • 断电不丢:即使你执行 kill -9,因为 timerwheel 和 timerlog 都在磁盘上,重启后,Broker 看一眼 checkpoint 文件,时针调回上一次断电的刻度,继续往前拨,没有一条定时消息会漏掉。